Compare commits
10 commits
e69d1c190f
...
52e657fb1e
| Author | SHA1 | Date | |
|---|---|---|---|
| 52e657fb1e | |||
| f64576372c | |||
| a15ea00c03 | |||
| 9d63e6482e | |||
| 09114f7c6a | |||
| d526001742 | |||
| bffd7faead | |||
| e3aad01baf | |||
| 5ccb2bfb52 | |||
| 0bafe9de01 |
18 changed files with 1357 additions and 16 deletions
11
.env.example
11
.env.example
|
|
@ -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
|
||||||
|
|
|
||||||
62
README.md
62
README.md
|
|
@ -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
72
app/api/schwab_auth.py
Normal 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"})
|
||||||
|
|
@ -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()]
|
||||||
|
|
|
||||||
|
|
@ -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
218
app/market/schwab.py
Normal 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()
|
||||||
|
|
@ -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)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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.
|
||||||
|
|
|
||||||
2
main.py
2
main.py
|
|
@ -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)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
255
scripts/check_schwab.py
Normal 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
98
scripts/check_stream.py
Normal 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))
|
||||||
|
|
@ -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();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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">
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
||||||
|
|
|
||||||
41
tests/test_schwab_callback.py
Normal file
41
tests/test_schwab_callback.py
Normal 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
253
tests/test_schwab_source.py
Normal 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]
|
||||||
Loading…
Reference in a new issue