Compare commits

..

10 commits

Author SHA1 Message Date
52e657fb1e Stream real-time /ES ticks so the candle 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, which reads as a dead feed.
LEVEL_ONE_FUTURES carries real trades on the same socket and the same login — no
extra REST call, no extra rate limit — and reports delayed: False on this
account. It was verified back in M6 and never subscribed to. It is now, building
a forming bar for the current minute that the authoritative CHART_FUTURES bar
then supersedes.

Three constraints shaped it, each a real bug avoided:

- Tick bars 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 again on every update. Runtime.on_bar returns early for
  an unclosed bar: store, set price, broadcast, stop.
- Emissions are throttled, SCHWAB_TICK_SECONDS default 1.0, because /ES trades
  many times a second and each emission is a store write plus a broadcast to
  every open socket. Negative drops the Level 1 subscription entirely.
- A tick for a minute CHART_FUTURES has already closed is dropped, or a late
  trade would overwrite a settled exchange bar with a partial one.

Bid-only updates are skipped rather than carried forward: a bid is not a trade
and must not extend a candle's high or low. Alerts stay on closed bars — a level
is judged on a settled bar, not a price that may not last the minute — which
needed no change, since on_bar already gated on closed.

Verified against the live socket: 15 forming bars and 2 closed bars in 100
seconds, the closed bar superseding each forming minute. Verified in a browser:
the last candle's high and low visibly extend within the minute, no console
errors. 85 tests pass, four of them new.

The plan gains the cold-restart options asked for: make seeding non-quadratic
first, then persist cooldowns, then persist bars.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 06:21:35 -05:00
f64576372c Add Font Awesome 7.3.1 and move the tool glyphs onto it
The sidebar drew its tool icons with box-drawing characters, which render at the
mercy of whatever font the system picks. Two are converted as the first users of
the icon set: a trend arrow for Trendline, a rule for Price level.

Served from unpkg, matching Vue and Lightweight Charts, and pinned for the same
reason they are — an unpinned icon set is a silent redesign on someone else's
release. cdnjs was the first choice and does not carry 7.3.1 on that path; it
404s, which in a browser shows up only as icons that quietly fail to their
fallback font at zero width. Verified rendering: computed family is
"Font Awesome 7 Free" at 17.5px with glyph content, no failed requests.

Lucide 1.31.0 is the lighter alternative if the weight ever matters — stroke SVGs
that can be inlined, dropping the CDN dependency entirely.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 06:10:28 -05:00
a15ea00c03 Keep level overlays on the candle grid so bars keep their width
The hourly chart drew 160 candles as unreadable slivers. Level points are
sampled on their own timeframe — VWAP every minute, the daily averages once a
session — and every distinct timestamp claims a slot on the chart's shared
scale. A minute-resolution VWAP on an hourly chart therefore spread 160 candles
across 908 slots, five to six times wider than the candles they belonged to:

  before   1h  inView 160  slots 908  ratio 5.46
  after    1h  inView 160  slots 161  ratio 1.01

Snapping ma and vwap points onto the candle grid fixes the density without
changing the line: points collapse onto the candle at or after them, and the
newest wins. Trendlines keep the raw path — lineData interpolates between two
anchors, and snapping those would move the geometry the user drew.

Two details worth keeping. Points past the final candle are clamped onto it
rather than 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 straight back on the scale. And the viewport is re-applied
once after the level series have loaded, because setBars runs before them and
the chart holds the width it derived from the previous timeframe's density —
consumed rather than reapplied, so the minutely VWAP resync cannot yank the view
back from wherever it has been panned.

The method is snapPointsToBars, not snapToBars: this.snapToBars already exists
as the "Snap to highs/lows" boolean, and the assignment silently replaced the
prototype method with true.

Verified in a browser across all six timeframes — 1m, 5m, 15m, 30m, 1h and 1d
each show 160 candles at ratio 0.99–1.01 with no console errors.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 06:07:44 -05:00
9d63e6482e Anchor the chart viewport by time, not by logical index
The chart sat about ten hours behind a perfectly healthy feed. The header price
updated live while the last candle stayed put, which reads as a dead stream.

Every layer checked out in isolation, because every layer was correct: /api/bars
served bars to the current minute from schwab, store.put keeps them strictly
ascending, the WebSocket snapshot delivered 1000 ascending bars ending at the
live edge, and the browser received all of it plus a bar event every minute.
Interrogating the page's own chart object is what separated "the data is
missing" from "the data is off-screen":

  seriesLen 1000  seriesLast 10:48 (7786.25)   data complete
  visible   08-07T20:41 -> 08-10T00:35         viewport 617 bars too far left

617 is exactly the daily bar count. setBars derived 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. A logical index addresses
the chart's shared time scale — the union of every series' time points — so
prepending those points renumbered every index and slid the view off the live
edge one tick after it had been set correctly. A time range names the instant
instead, and later series cannot move it.

Worth keeping: screenshots alone were actively misleading. The stale time axis
showed a Friday-to-Sunday gap that read as an ordinary session break, so the
view looked plausible while being ten hours wrong.

The plan gains a dated session log (§16) for this and for two things that cost a
detour today — the image needing a rebuild for schwab-py, which the bind mount
hides, and the ~82 second startup during which the port refuses connections and
an open tab logs a wall of ERR_CONNECTION_REFUSED. That slow seed is recorded,
not fixed: it replays every bar through on_bar and rebuilds all five MA levels
per daily bar.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 05:53:13 -05:00
09114f7c6a Switch the live feed to Schwab and document the rate limits
LIVE_SOURCE=schwab locally. Seeding stays on Yahoo, which is enforced in the
factory rather than left to configuration.

The portal's "Order Limit: 120" caps orders per minute and this app places none;
Schwab's separate REST limit is commonly cited at the same number. Neither binds
here, because streaming is not REST — one socket, bars pushed, essentially no
REST traffic in steady state. Worth recording as a reason the stream beats the
polling fallback beyond latency: polling quotes once a second would have sat at
half the limit permanently.

Reconnects are the exception. Each calls get_user_preferences() for the socket
URL, and the retry backoff is five seconds, so a sustained outage costs about
twelve REST calls a minute — under the limit, but a reason not to shorten that
backoff without thinking.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 05:26:47 -05:00
d526001742 Add the Schwab live source: real-time /ES minute bars
Verified against a live account before and after writing it. CHART_FUTURES
delivers one true-OHLCV minute bar per symbol per minute, LEVEL_ONE_FUTURES
reports delayed: false, and consecutive bars arrived sixty seconds apart through
the production code path.

Yahoo stays. Schwab serves no futures history whatever, so seed_source resolves
to Yahoo even when SEED_SOURCE=schwab is asked for — the pairing is the intended
configuration rather than a fallback. The symbols differ, ES=F against /ES, so
Settings.live_symbol picks the live one while seeding always uses Yahoo's.

Three findings worth keeping, each of which cost a round trip:

- get_quote() singular returns the wrong instrument entirely. It puts the symbol
  in the URL path, where the leading slash is normalised away, so /ES resolves to
  Eversource Energy at $72 and returns HTTP 200 with a populated body. Only
  get_quotes() plural, which passes symbols as a query parameter, returns the
  future. A 200 is not evidence; assetMainType is.
- Streaming requires the Accounts and Trading product. StreamClient.login() reads
  /trader/v1/userPreference for its socket URL, and that path does not exist in
  Market Data Production.
- /ES resolves to the active contract on Schwab's side, so the contract roll
  handling the plan left open needs no code.

The stream drops the oldest queued message rather than stalling the socket, and
surfaces a dead pump task instead of waiting forever on a queue nothing fills.
schwab-py moves into requirements.txt, imported only when LIVE_SOURCE=schwab.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 05:23:20 -05:00
bffd7faead Validate Schwab credentials before sending anyone to a browser
The first attempt failed with invalid_client and no way to tell why: the key,
the secret, the callback, or an app not yet propagated all look identical from
the browser, which shows raw JSON. It turned out to be a key clipped by one
character on paste.

Two preflight checks now say which. The authorize endpoint is asked whether it
recognises the key. The token endpoint is asked to exchange a deliberately
invalid code, which separates bad credentials from a bad grant — it
authenticates the key and secret over HTTP Basic before it looks at the code, so
invalid_client means the pair is wrong and invalid_grant means the pair is fine.

That second check matters more than it sounds. The secret is not used at all
during login, so a truncated one survives the whole browser round trip and only
surfaces at the exchange, by which point the authorisation code has been spent
and the flow has to start over.

Neither check can tell a wrong value from an app that is not live yet — an
invented key produces the identical response, verified — and both say so rather
than guessing.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 05:02:04 -05:00
e3aad01baf Guard that the OAuth callback stays reachable under CHART_AUTH_TOKEN
The callback is registered with the provider and has to answer an unauthenticated
browser redirect. It is exempt by construction — a separate router without the
token dependency — but nothing held that in place, and the failure would only
appear in production, where the token is the one setting that differs from
local, at the last step of a login flow.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 04:50:12 -05:00
5ccb2bfb52 Add Schwab settings, credential placeholders and an entitlement probe
config.py had no Schwab fields at all, so the keys listed in .env.example were
being silently dropped by extra="ignore". They exist now, blank, and nothing
reads them while live_source is yahoo.

The token path moves under data/, which is the Coolify persistent volume. Left
at the repository root it would vanish on every rebuild, and re-authenticating
is an interactive browser flow, not something a deploy can do for itself.

scripts/check_schwab.py answers empirically what the app is entitled to rather
than inferring it from documentation: whether the credentials authenticate,
whether /ES quotes return (futures market data is a separate entitlement from
equities), and whether the streamer bootstrap responds.

That last one is the decision. StreamClient.login() reads
/trader/v1/userPreference for its socket URL and credentials, and that path
belongs to the Accounts and Trading product — so an app registered for Market
Data Production alone cannot stream, and CHART_FUTURES is unreachable until the
app adds it. The script reports which of the two paths is open instead of
leaving it to be discovered halfway through an implementation.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 04:48:29 -05:00
0bafe9de01 Add the OAuth callback endpoint at /api/qt
Schwab requires an HTTPS callback. The usual answer is https://127.0.0.1:8182
behind 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 take the redirect
itself.

Unauthenticated by necessity: the provider redirects a browser here and cannot
attach the chart token, so it sits alongside /health and /version. It is inert —
nothing is stored, and the page echoes only the query string of the request that
produced it, which the caller already has in their address bar. Retaining the
code would let a later anonymous visitor read it.

The path and the page are both deliberately unrevealing. That is not a security
control; it just avoids advertising which brokerage this host talks to. Treat
the path as fixed — changing a registered callback means editing the app, which
can send it back through approval.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 04:45:31 -05:00
18 changed files with 1357 additions and 16 deletions

View file

@ -6,6 +6,17 @@ YAHOO_POLL_SECONDS=20
SEED_1H_RANGE=730d SEED_1H_RANGE=730d
SEED_1M_RANGE=8d 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 # Chart and analysis
TIMEFRAMES=1m,5m,15m,30m,1h,1d TIMEFRAMES=1m,5m,15m,30m,1h,1d
BASE_TIMEFRAMES=1m,30m,1d BASE_TIMEFRAMES=1m,30m,1d

View file

@ -172,6 +172,68 @@ knowing:
If it lands near other levels it clusters with them and the score adds up, so a typed 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. 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 and ntfy
Alerts are evaluated **server-side**, once per closed 1m bar, by a single engine Alerts are evaluated **server-side**, once per closed 1m bar, by a single engine

72
app/api/schwab_auth.py Normal file
View file

@ -0,0 +1,72 @@
"""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,6 +32,20 @@ class Settings(BaseSettings):
ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200" ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200"
daily_anchor_et: str = "18:00" daily_anchor_et: str = "18:00"
manual_lines_path: Path = Path("./data/manual_lines.json") 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 confluence_min_score: float = 28
# Four hours, chosen from the sweep in scripts/calibrate_alerts.py. Suppression is # 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 # per price zone, so an unrelated zone still alerts immediately; this only
@ -42,6 +56,15 @@ class Settings(BaseSettings):
chart_auth_token: str = "" chart_auth_token: str = ""
replay_file: Path | None = None 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 @property
def enabled_timeframes(self) -> list[Timeframe]: def enabled_timeframes(self) -> list[Timeframe]:
return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()] return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()]

View file

@ -7,15 +7,27 @@ from app.market.yahoo import YahooSource
def live_source(settings: Settings) -> MarketDataSource: def live_source(settings: Settings) -> MarketDataSource:
if settings.live_source == "yahoo": if settings.live_source == "yahoo":
return YahooSource(settings.yahoo_poll_seconds) 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: if settings.live_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file) return ReplaySource(settings.replay_file)
raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}") raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}")
def seed_source(settings: Settings) -> MarketDataSource | None: 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": if settings.seed_source == "none":
return None return None
if settings.seed_source == "yahoo": if settings.seed_source in ("yahoo", "schwab"):
return YahooSource(settings.yahoo_poll_seconds) return YahooSource(settings.yahoo_poll_seconds)
if settings.seed_source == "replay" and settings.replay_file: if settings.seed_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file) return ReplaySource(settings.replay_file)

218
app/market/schwab.py Normal file
View file

@ -0,0 +1,218 @@
"""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,11 +22,17 @@ class StreamService:
self._handlers.append(handler) self._handlers.append(handler)
async def seed( async def seed(
self, source: MarketDataSource | None, tf: Timeframe, range_: str self,
source: MarketDataSource | None,
tf: Timeframe,
range_: str,
symbol: str | None = None,
) -> None: ) -> 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(): if source is None or not source.supports_history():
return return
bars = await source.history(self.symbol, tf, None, None, range_=range_) bars = await source.history(symbol or self.symbol, tf, None, None, range_=range_)
for bar in bars: for bar in bars:
await self._emit(bar) await self._emit(bar)

View file

@ -50,10 +50,24 @@ class Runtime:
self.settings.confluence_min_score, self.settings.alert_cooldown_seconds self.settings.confluence_min_score, self.settings.alert_cooldown_seconds
) )
self.levels = self.manual_lines.levels() self.levels = self.manual_lines.levels()
self.stream = StreamService(live_source(self.settings), self.settings.yahoo_symbol) self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
self.stream.add_handler(self.on_bar) self.stream.add_handler(self.on_bar)
async def on_bar(self, bar: Bar) -> None: 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 evaluate_alerts = False
for aggregated in self.aggregator.update(bar): for aggregated in self.aggregator.update(bar):
self.store.put(aggregated) self.store.put(aggregated)
@ -181,8 +195,14 @@ class Runtime:
async def start(self) -> asyncio.Task: async def start(self) -> asyncio.Task:
try: try:
source = seed_source(self.settings) source = seed_source(self.settings)
await self.stream.seed(source, Timeframe.H1, self.settings.seed_1h_range) # Always the Yahoo symbol: Schwab has no history to seed from.
await self.stream.seed(source, Timeframe.M1, self.settings.seed_1m_range) 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
)
except Exception: except Exception:
# A transient seed failure must not prevent the live stream or UI starting. # A transient seed failure must not prevent the live stream or UI starting.
pass pass

View file

@ -168,10 +168,37 @@ 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. Yahoo's native 1d bars may be used *only* for multi-year context, clearly labelled.
## 2.2 Verify these when adding the Schwab source (M6, not before) ## 2.2 Schwab entitlements — **answered empirically 2026-08-10**
These are load-bearing unknowns for Schwab specifically. They no longer block the All verified against a live production app holding both Market Data Production and
project — M0–M5 run entirely on Yahoo. If any fails, stop and report. 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
1. **Futures market-data entitlement.** It is not publicly documented whether 1. **Futures market-data entitlement.** It is not publicly documented whether
`CHART_FUTURES` requires futures trading approval or a CME non-professional market `CHART_FUTURES` requires futures trading approval or a CME non-professional market
@ -783,6 +810,26 @@ const chartApi = shallowRef(null); // ✅
// const chart = ref(null); // ❌ will appear to work, then misbehave // 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: Structure:
- `chart.js` — a plain, framework-free wrapper class owning the LWC instance: - `chart.js` — a plain, framework-free wrapper class owning the LWC instance:
`create(el)`, `setBars()`, `updateBar()`, `syncLevels(levels)`, `destroy()`. `create(el)`, `setBars()`, `updateBar()`, `syncLevels(levels)`, `destroy()`.
@ -1144,3 +1191,134 @@ enforces for data sources.
| Schwab token expiry (7 days) | Stream dies | Surface prominently in status bar; document re-auth | | 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 | | Vue reactivity wrapping chart objects | Perf collapse, odd bugs | `shallowRef`/`markRaw` — §9 |
| LWC v4 tutorials copied | Code silently wrong for v5 | `addSeries(SeriesType, ...)` only | | 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,6 +12,7 @@ from fastapi.staticfiles import StaticFiles
from app.api.meta import router as meta_router from app.api.meta import router as meta_router
from app.api.routes import router as api_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.api.ws import router as ws_router
from app.config import Settings from app.config import Settings
from app.runtime import Runtime from app.runtime import Runtime
@ -37,6 +38,7 @@ app = FastAPI(title="chart", lifespan=lifespan)
app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static")
app.include_router(meta_router) app.include_router(meta_router)
app.include_router(schwab_auth_router)
app.include_router(api_router) app.include_router(api_router)
app.include_router(ws_router) app.include_router(ws_router)

View file

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

255
scripts/check_schwab.py Normal file
View file

@ -0,0 +1,255 @@
"""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()

98
scripts/check_stream.py Normal file
View file

@ -0,0 +1,98 @@
"""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,6 +32,7 @@ class ConfluenceChart {
this.toolDownListener = null; this.toolDownListener = null;
this.toolMoveListener = null; this.toolMoveListener = null;
this.toolUpListener = null; this.toolUpListener = null;
this.pendingView = null;
} }
static TICK = 0.25; static TICK = 0.25;
@ -144,10 +145,22 @@ class ConfluenceChart {
setBars(bars) { setBars(bars) {
this.bars = bars; this.bars = bars;
this.candles.setData(bars.map(this.toCandle)); this.candles.setData(bars.map(this.toCandle));
this.chart.timeScale().setVisibleLogicalRange({ // Anchored by time, not by logical index. A logical index addresses the
from: Math.max(0, bars.length - 160), // chart's *shared* scale — the union of every series' time points — not
to: bars.length + 5, // 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);
}
requestAnimationFrame(() => this.renderAnchorHandles()); requestAnimationFrame(() => this.renderAnchorHandles());
} }
@ -190,6 +203,42 @@ 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) { syncLevels(levels) {
this.levels = levels; this.levels = levels;
this.syncPriceLines(levels); this.syncPriceLines(levels);
@ -227,7 +276,7 @@ class ConfluenceChart {
} }
let data; let data;
if (hasPoints) { if (hasPoints) {
data = (level.points || []).map(([time, value]) => ({ time, value })); data = this.snapPointsToBars((level.points || []).map(([time, value]) => ({ time, value })));
const latestTime = this.bars[this.bars.length - 1]?.t; const latestTime = this.bars[this.bars.length - 1]?.t;
const latestValue = data[data.length - 1]?.value; const latestValue = data[data.length - 1]?.value;
if (latestTime != null && latestValue != null && latestTime > data[data.length - 1].time) { if (latestTime != null && latestValue != null && latestTime > data[data.length - 1].time) {
@ -238,6 +287,16 @@ class ConfluenceChart {
} }
entry.series.setData(data); 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(); this.renderAnchorHandles();
} }

View file

@ -5,6 +5,10 @@
<meta name="viewport" content="width=device-width, initial-scale=1"> <meta name="viewport" content="width=device-width, initial-scale=1">
<title>/ES Confluence</title> <title>/ES Confluence</title>
<link rel="stylesheet" href="/static/style.css"> <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/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> <script src="https://unpkg.com/lightweight-charts@5.2.0/dist/lightweight-charts.standalone.production.js"></script>
</head> </head>
@ -34,7 +38,7 @@
<h2>Tools</h2> <h2>Tools</h2>
<div class="tool" :class="{armed: armedTool === 'trendline'}"> <div class="tool" :class="{armed: armedTool === 'trendline'}">
<button class="tool-head" @click="armTool('trendline')" :aria-pressed="armedTool === 'trendline'"> <button class="tool-head" @click="armTool('trendline')" :aria-pressed="armedTool === 'trendline'">
<span class="tool-glyph">╱</span>Trendline <span class="tool-glyph"><i class="fa-solid fa-arrow-trend-up"></i></span>Trendline
<span class="tool-state">{{ armedTool === 'trendline' ? 'drag on chart' : '' }}</span> <span class="tool-state">{{ armedTool === 'trendline' ? 'drag on chart' : '' }}</span>
</button> </button>
<div class="tool-body"> <div class="tool-body">
@ -51,7 +55,7 @@
</div> </div>
<div class="tool" :class="{armed: armedTool === 'level'}"> <div class="tool" :class="{armed: armedTool === 'level'}">
<button class="tool-head" @click="armTool('level')" :aria-pressed="armedTool === 'level'"> <button class="tool-head" @click="armTool('level')" :aria-pressed="armedTool === 'level'">
<span class="tool-glyph">─</span>Price level <span class="tool-glyph"><i class="fa-solid fa-minus"></i></span>Price level
<span class="tool-state">{{ armedTool === 'level' ? 'click on chart' : '' }}</span> <span class="tool-state">{{ armedTool === 'level' ? 'click on chart' : '' }}</span>
</button> </button>
<div class="tool-body"> <div class="tool-body">

View file

@ -5,6 +5,7 @@ from starlette.websockets import WebSocketDisconnect
from app.api.meta import router as meta_router from app.api.meta import router as meta_router
from app.api.routes import router as api_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.api.ws import router as ws_router
from app.config import Settings from app.config import Settings
from app.runtime import Runtime from app.runtime import Runtime
@ -19,6 +20,7 @@ def client(tmp_path):
) )
app = FastAPI() app = FastAPI()
app.include_router(meta_router) app.include_router(meta_router)
app.include_router(schwab_auth_router)
app.include_router(api_router) app.include_router(api_router)
app.include_router(ws_router) app.include_router(ws_router)
app.state.runtime = Runtime(settings) app.state.runtime = Runtime(settings)
@ -31,6 +33,29 @@ def test_open_when_no_token_configured(client):
assert client("").get("/api/bars").status_code == 200 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): def test_rejects_missing_token(client):
assert client("s3cret").get("/api/bars").status_code == 401 assert client("s3cret").get("/api/bars").status_code == 401

View file

@ -0,0 +1,41 @@
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

253
tests/test_schwab_source.py Normal file
View file

@ -0,0 +1,253 @@
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]