chart/app/runtime.py
Chris Amow 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

195 lines
8.4 KiB
Python

import asyncio
import logging
from dataclasses import dataclass, field
from app.analysis.alerts import Alert, AlertEngine
from app.bars.models import Bar, Timeframe
from app.bars.aggregator import Aggregator
from app.analysis.bar_space import price_in_bar_space
from app.analysis.horizontals import build_prior_day_levels
from app.analysis.levels import Level
from app.analysis.moving_averages import build_ma_levels
from app.analysis.vwap import build_vwap_level
from app.analysis.confluence import Cluster, cluster_levels
from app.analysis.indicators import atr
from app.analysis.manual_lines import ManualLineStore
from app.bars.store import InMemoryBarStore
from app.config import Settings
from app.market.factory import live_source, seed_source
from app.market.stream import StreamService
from app.notify.ntfy import send_ntfy
logger = logging.getLogger(__name__)
@dataclass
class Runtime:
settings: Settings
store: InMemoryBarStore = field(init=False)
stream: StreamService = field(init=False)
subscribers: set[asyncio.Queue[dict]] = field(default_factory=set)
aggregator: Aggregator = field(init=False)
levels: list[Level] = field(default_factory=list)
clusters: list[Cluster] = field(default_factory=list)
price: float | None = None
atr15: float = 0.0
manual_lines: ManualLineStore = field(init=False)
ma_levels: list[Level] = field(default_factory=list)
alert_engine: AlertEngine = field(init=False)
_sent_levels: dict[str, dict] = field(default_factory=dict)
_notify_tasks: set[asyncio.Task] = field(default_factory=set)
def __post_init__(self) -> None:
self.store = InMemoryBarStore(self.settings.max_bars_per_tf)
self.aggregator = Aggregator(self.settings.enabled_timeframes)
self.manual_lines = ManualLineStore(self.settings.manual_lines_path)
# One engine for the process, not one per browser connection. Cooldowns
# are only meaningful if they outlive a page reload, and a phone push
# must not depend on a tab being open to produce it.
self.alert_engine = AlertEngine(
self.settings.confluence_min_score, self.settings.alert_cooldown_seconds
)
self.levels = self.manual_lines.levels()
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
self.stream.add_handler(self.on_bar)
async def on_bar(self, bar: Bar) -> None:
evaluate_alerts = False
for aggregated in self.aggregator.update(bar):
self.store.put(aggregated)
self.broadcast({"type": "bar", "bar": aggregated})
if self.settings.ma_sets.get(aggregated.tf):
self.rebuild_levels()
if aggregated.tf is Timeframe.M1 and aggregated.closed:
self.price = aggregated.c
evaluate_alerts = True
if evaluate_alerts:
values = atr(self.store.get(Timeframe.M15), 14)
self.atr15 = next((value for value in reversed(values) if value is not None), 0.0)
# VWAP re-prices every minute, so levels are rebuilt here too. The
# broadcast is a delta, which is what keeps that affordable.
self.rebuild_levels()
self.rebuild_clusters(evaluate_alerts=True)
def broadcast(self, event: dict) -> None:
for queue in self.subscribers.copy():
if queue.full():
queue.get_nowait()
queue.put_nowait(event)
def rebuild_levels(self) -> None:
self.ma_levels = build_ma_levels(
{tf: self.store.get(tf) for tf in self.settings.ma_sets},
self.settings.ma_sets,
)
minute_bars = self.store.get(Timeframe.M1)
manual = self.manual_lines.levels()
self.position_manual_levels(manual, minute_bars)
self.levels = (
self.ma_levels
+ build_prior_day_levels(self.store.get(Timeframe.D1), self.price)
+ build_vwap_level(minute_bars)
+ manual
)
self.broadcast_level_delta()
self.rebuild_clusters()
@staticmethod
def position_manual_levels(levels: list[Level], bars: list[Bar]) -> None:
"""Price sloped lines across bars rather than seconds.
The chart spaces bars evenly, so the line a person drew advances per bar.
Pricing it per second instead put the alert somewhere the line visibly
was not — 147 points out across a weekend.
"""
if not bars:
return
times = [bar.t for bar in bars]
now = times[-1]
for level in levels:
if level.slope:
level.current_p = price_in_bar_space(level, times, now)
def broadcast_level_delta(self) -> None:
"""Send only levels whose serialised form actually changed.
A daily moving average carries hundreds of points and changes once a
session; VWAP changes every minute. Broadcasting the whole set on the
VWAP cadence would push the entire history every minute, so subscribers
get a delta and merge it by id.
"""
current = {level.id: level.to_dict() for level in self.levels}
changed = [value for id_, value in current.items() if self._sent_levels.get(id_) != value]
removed = [id_ for id_ in self._sent_levels if id_ not in current]
self._sent_levels = current
if changed or removed:
self.broadcast({"type": "levels", "changed": changed, "removed": removed})
def rebuild_clusters(self, evaluate_alerts: bool = False) -> None:
if self.price is None or self.stream.last_bar_t is None:
return
self.clusters = cluster_levels(self.levels, self.stream.last_bar_t, self.price, self.atr15)
self.broadcast({"type": "clusters", "price": self.price, "clusters": self.clusters})
if evaluate_alerts:
# Evaluated over every level, deliberately ignoring per-connection
# layer preferences: those are a display choice made in one browser,
# and a push notification has no business depending on them.
self.dispatch_alerts(
self.alert_engine.evaluate(
self.clusters,
self.price,
self.atr15,
self.stream.last_bar_t,
self.stream.symbol,
)
)
def dispatch_alerts(self, alerts: list[Alert]) -> None:
tripped: set[str] = set()
for alert in alerts:
self.broadcast({"type": "alert", "cluster": alert.cluster, "message": alert.message})
task = asyncio.create_task(self.notify(alert.message))
# Held so the task is not garbage collected mid-flight.
self._notify_tasks.add(task)
task.add_done_callback(self._notify_tasks.discard)
tripped.update(alert.tripped)
if tripped:
self.disarm(tripped)
def disarm(self, ids: set[str]) -> None:
"""A hand-placed level fires once, then waits to be re-armed."""
changed = False
for line_id in ids:
try:
self.manual_lines.update(line_id, {"armed": False})
changed = True
except KeyError:
continue
# Rebuilt after the loop, not inside it: rebuild_levels re-enters
# rebuild_clusters, and doing that mid-dispatch would rewrite the very
# clusters being iterated.
if changed:
self.rebuild_levels()
async def notify(self, message: str) -> None:
try:
await send_ntfy(self.settings.ntfy_server, self.settings.ntfy_topic, message)
except Exception:
# A push outage must not take down the stream or the sockets.
logger.warning("ntfy delivery failed", exc_info=True)
async def start(self) -> asyncio.Task:
try:
source = seed_source(self.settings)
# Always the Yahoo symbol: Schwab has no history to seed from.
seed_symbol = self.settings.yahoo_symbol
await self.stream.seed(
source, Timeframe.H1, self.settings.seed_1h_range, seed_symbol
)
await self.stream.seed(
source, Timeframe.M1, self.settings.seed_1m_range, seed_symbol
)
except Exception:
# A transient seed failure must not prevent the live stream or UI starting.
pass
return asyncio.create_task(self.stream.run(), name="market-stream")