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

209 lines
9.1 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:
# A tick-built bar is provisional and arrives many times a minute. It
# updates the last candle and the live price, and stops there.
#
# It must not reach the aggregator: that accumulates volume with
# `current.v += incoming.v`, so re-sending the same forming minute would
# add its volume to every higher timeframe again on each update. Alerts
# stay on closed bars for the same reason they always were — a level is
# judged on a settled bar, not on a price that may not last the minute.
if not bar.closed:
self.store.put(bar)
self.price = bar.c
self.broadcast({"type": "bar", "bar": bar})
return
evaluate_alerts = False
for aggregated in self.aggregator.update(bar):
self.store.put(aggregated)
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")