449 lines
20 KiB
Python
449 lines
20 KiB
Python
import asyncio
|
|
import logging
|
|
import time
|
|
from dataclasses import dataclass, field, replace
|
|
from typing import ClassVar
|
|
|
|
from app.analysis.alerts import Alert, AlertEngine
|
|
from app.analysis.event_log import EventLog
|
|
from app.bars.models import Bar, Timeframe
|
|
from app.bars.aggregator import Aggregator
|
|
from app.bars.session import bucket_start
|
|
from app.analysis.bar_space import price_in_bar_space, price_in_timeframe_space
|
|
from app.analysis.horizontals import build_prior_day_levels
|
|
from app.analysis.levels import Level, LevelKind
|
|
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.analysis.user_prefs import UserPrefStore
|
|
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)
|
|
user_prefs: UserPrefStore = 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)
|
|
# The loop that owns the subscriber queues. Set once the app is running;
|
|
# None while a test drives the runtime directly.
|
|
_loop: asyncio.AbstractEventLoop | None = None
|
|
# True while history is being replayed, so on_bar accumulates quietly.
|
|
seeding: bool = False
|
|
# Coalescing window for level rebuilds; see request_rebuild.
|
|
REBUILD_INTERVAL: ClassVar[float] = 0.25
|
|
_last_rebuild: float = 0.0
|
|
_rebuild_pending: bool = False
|
|
# Seconds the loop ran late, worst since start and most recent sample.
|
|
loop_lag_worst: float = 0.0
|
|
loop_lag_recent: float = 0.0
|
|
_lag_task: asyncio.Task | None = None
|
|
_rebuild_task: asyncio.Task | None = None
|
|
_token_task: asyncio.Task | None = None
|
|
schwab_login: object | None = None
|
|
events: EventLog = field(init=False)
|
|
|
|
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)
|
|
self.user_prefs = UserPrefStore(self.settings.user_prefs_path)
|
|
self.events = EventLog(self.settings.events_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.settings.alert_state_path,
|
|
self.settings.alert_timezone,
|
|
)
|
|
self.levels = self.manual_lines.levels()
|
|
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
|
|
self.stream.add_handler(self.on_bar)
|
|
self.stream.on_drop = self._on_stream_drop
|
|
|
|
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})
|
|
for provisional in self.provisional_higher(bar):
|
|
self.store.put(provisional)
|
|
self.broadcast({"type": "bar", "bar": provisional})
|
|
return
|
|
|
|
evaluate_alerts = False
|
|
for aggregated in self.aggregator.update(bar):
|
|
self.store.put(aggregated)
|
|
if self.seeding:
|
|
# Seeding replays years of history through this method. Every
|
|
# bar used to broadcast to nobody and rebuild every level from
|
|
# scratch, which is the whole of the startup stall — 19.5
|
|
# seconds of measured loop lag and two minutes before the port
|
|
# opened. Bars are accumulated here and the levels are built
|
|
# once, at the end, from the finished store.
|
|
continue
|
|
self.broadcast({"type": "bar", "bar": aggregated})
|
|
if self.settings.ma_sets.get(aggregated.tf):
|
|
self.request_rebuild()
|
|
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.request_rebuild()
|
|
self.rebuild_clusters(evaluate_alerts=True)
|
|
|
|
def provisional_higher(self, bar: Bar) -> list[Bar]:
|
|
"""Higher-timeframe bars including the minute still being traded.
|
|
|
|
The aggregator cannot be asked for these: it accumulates with
|
|
``current.v += incoming.v``, so re-feeding the same forming minute on
|
|
every tick would add its volume to each higher timeframe again and
|
|
again. Its committed state — every minute that has actually closed — is
|
|
combined with the live minute here instead, without mutating it. The
|
|
next closed minute goes through the aggregator normally and replaces
|
|
what this produced, because the store keys on the bucket's timestamp.
|
|
"""
|
|
out: list[Bar] = []
|
|
for tf in self.settings.enabled_timeframes:
|
|
if tf is Timeframe.M1:
|
|
continue
|
|
start = bucket_start(bar.t, tf)
|
|
base = self.aggregator.forming.get(tf)
|
|
if base is None or start > base.t:
|
|
# The live minute opens a bucket the aggregator has not started.
|
|
out.append(replace(bar, tf=tf, t=start, closed=False))
|
|
continue
|
|
if start < base.t:
|
|
continue
|
|
out.append(
|
|
replace(
|
|
base,
|
|
h=max(base.h, bar.h),
|
|
l=min(base.l, bar.l),
|
|
c=bar.c,
|
|
v=base.v + bar.v,
|
|
closed=False,
|
|
)
|
|
)
|
|
return out
|
|
|
|
def broadcast(self, event: dict) -> None:
|
|
"""Publish an event to every subscriber, from any thread.
|
|
|
|
asyncio.Queue is not thread-safe: it wakes a waiting consumer by
|
|
resolving a Future, which only the loop thread may do. The mutating
|
|
routes are sync `def`, so FastAPI runs them in a threadpool, and they
|
|
reach here through rebuild_levels — writing the queue directly from
|
|
there can drop a socket's wakeup. The visible symptom is a drawing made
|
|
in one browser not reaching another until the next market tick, which
|
|
is why it has gone unnoticed: the stream ticks about once a second and
|
|
covers it over.
|
|
"""
|
|
loop = self._loop
|
|
if loop is None or self._on_loop_thread(loop):
|
|
self._publish(event)
|
|
return
|
|
loop.call_soon_threadsafe(self._publish, event)
|
|
|
|
@staticmethod
|
|
def _on_loop_thread(loop: asyncio.AbstractEventLoop) -> bool:
|
|
try:
|
|
return asyncio.get_running_loop() is loop
|
|
except RuntimeError:
|
|
# No loop in this thread at all, so certainly not that one.
|
|
return False
|
|
|
|
def _publish(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()
|
|
|
|
def position_manual_levels(self, 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
|
|
minute_times = [bar.t for bar in bars]
|
|
now = minute_times[-1]
|
|
for level in levels:
|
|
if not level.slope:
|
|
continue
|
|
if not self.settings.trendline_source_geometry:
|
|
level.current_p = price_in_bar_space(level, minute_times, now)
|
|
continue
|
|
source_times = [bar.t for bar in self.store.get(level.tf)]
|
|
level.current_p = price_in_timeframe_space(level, source_times, level.tf, now)
|
|
level.geometry_resolved = level.current_p is not None
|
|
|
|
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,
|
|
self.watched_ma_levels(),
|
|
)
|
|
)
|
|
|
|
def ma_alerts(self) -> dict:
|
|
stored = self.user_prefs.get("ma_alerts") or {}
|
|
return {"1d": list(stored.get("1d") or [])}
|
|
|
|
def set_ma_alerts(self, value: dict) -> dict:
|
|
self.user_prefs.put("ma_alerts", value)
|
|
return self.ma_alerts()
|
|
|
|
def watched_ma_levels(self) -> list[Level]:
|
|
armed = set(self.ma_alerts().get("1d", []))
|
|
return [
|
|
level for level in self.levels
|
|
if level.kind is LevelKind.MA
|
|
and level.tf is Timeframe.D1
|
|
and level.period in armed
|
|
]
|
|
|
|
def _on_stream_drop(self, error: str) -> None:
|
|
kind = "auth" if "invalid_grant" in error or "Refresh token" in error else "stream"
|
|
self.events.add(kind, error.split("\n", 1)[0][:200])
|
|
|
|
def dispatch_alerts(self, alerts: list[Alert]) -> None:
|
|
tripped: set[str] = set()
|
|
for alert in alerts:
|
|
self.events.add("alert", alert.message, number=alert.number, at=alert.at)
|
|
self.broadcast({
|
|
"type": "alert",
|
|
"cluster": alert.cluster,
|
|
"message": alert.message,
|
|
# Same number the push carries, so a phone and a screen agree.
|
|
"number": alert.number,
|
|
"at": alert.at,
|
|
})
|
|
task = asyncio.create_task(self.notify(alert.push or 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)
|
|
|
|
def request_rebuild(self) -> None:
|
|
"""Rebuild levels soon, at most REBUILD_INTERVAL apart.
|
|
|
|
A rebuild costs a pass over every moving average and a diff of their
|
|
points, so doing one per bar is fine at one bar a minute and ruinous in
|
|
a burst. Yahoo's first poll emits a whole day of minutes at once, which
|
|
measured as twenty seconds of loop lag. Coalescing turns that into a
|
|
handful of rebuilds without changing what subscribers eventually see:
|
|
the pending flag guarantees a trailing rebuild, so the last bar of a
|
|
burst is never the one that gets dropped.
|
|
"""
|
|
now = time.monotonic()
|
|
if now - self._last_rebuild >= self.REBUILD_INTERVAL:
|
|
self._last_rebuild = now
|
|
self._rebuild_pending = False
|
|
self.rebuild_levels()
|
|
return
|
|
self._rebuild_pending = True
|
|
|
|
async def rebuild_watch(self, interval: float = 0.25) -> None:
|
|
"""Flush a rebuild that request_rebuild deferred during a burst."""
|
|
while True:
|
|
await asyncio.sleep(interval)
|
|
if self._rebuild_pending:
|
|
self._rebuild_pending = False
|
|
self._last_rebuild = time.monotonic()
|
|
self.rebuild_levels()
|
|
|
|
def settle_after_seed(self) -> None:
|
|
"""Derive everything the replay deliberately skipped, once.
|
|
|
|
on_bar normally maintains price, ATR and the level set as each bar
|
|
arrives. During a seed it only fills the store, so the same state is
|
|
computed here from the finished history — one pass instead of one per
|
|
bar. Alerts are not evaluated: a level touched two years ago is not news,
|
|
and firing on replayed history is how a deploy used to re-alert.
|
|
"""
|
|
minute_bars = self.store.get(Timeframe.M1)
|
|
if minute_bars:
|
|
self.price = minute_bars[-1].c
|
|
values = atr(self.store.get(Timeframe.M15), 14)
|
|
self.atr15 = next((value for value in reversed(values) if value is not None), 0.0)
|
|
self.rebuild_levels()
|
|
|
|
async def loop_lag_watch(self, interval: float = 0.1) -> None:
|
|
"""Measure how late the event loop is running its own timers.
|
|
|
|
The loop is single-threaded and everything shares it: the market
|
|
stream, every WebSocket, and any CPU work that has strayed onto it.
|
|
When something blocks, the symptom reaching a person is "the chart
|
|
feels laggy" — unfalsifiable. This turns it into a number.
|
|
|
|
Scheduling drift is the honest measure: sleep for a known interval and
|
|
see how much longer it actually took.
|
|
"""
|
|
while True:
|
|
before = time.perf_counter()
|
|
await asyncio.sleep(interval)
|
|
lag = (time.perf_counter() - before) - interval
|
|
if lag > self.loop_lag_worst:
|
|
self.loop_lag_worst = lag
|
|
self.loop_lag_recent = lag
|
|
|
|
async def start(self) -> asyncio.Task:
|
|
# Captured here so a threadpool route can post events back to the loop
|
|
# that owns the queues, rather than touching them across threads.
|
|
self._loop = asyncio.get_running_loop()
|
|
self.seeding = True
|
|
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
|
|
)
|
|
# One-hour history cannot reconstruct a 30-minute candle. Yahoo
|
|
# supplies roughly 60 days natively; the newer overlap is replaced
|
|
# below by bars aggregated from the finer 1m seed.
|
|
await self.stream.seed(
|
|
source, Timeframe.M30, self.settings.seed_30m_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
|
|
finally:
|
|
self.seeding = False
|
|
self.settle_after_seed()
|
|
self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag")
|
|
self._rebuild_task = asyncio.create_task(self.rebuild_watch(), name="level-rebuild")
|
|
keep_alive = getattr(self.stream.source, "keep_alive", None)
|
|
if keep_alive is not None:
|
|
self._token_task = asyncio.create_task(keep_alive(), name="token-keepalive")
|
|
return asyncio.create_task(self.stream.run(), name="market-stream")
|
|
|
|
def needs_login(self) -> bool:
|
|
error = self.stream.last_error or ""
|
|
return "invalid_grant" in error or "Refresh token" in error
|
|
|
|
def start_schwab_login(self) -> str:
|
|
from app.api.schwab_auth import begin_login
|
|
|
|
context = begin_login(self.settings)
|
|
self.schwab_login = context
|
|
return context.authorization_url
|
|
|
|
def finish_schwab_login(self, redirect_url: str) -> None:
|
|
from app.api.schwab_auth import complete_login
|
|
|
|
context = self.schwab_login
|
|
if context is None:
|
|
raise RuntimeError("No login in progress")
|
|
complete_login(self.settings, context, redirect_url)
|
|
self.schwab_login = None
|
|
reset = getattr(self.stream.source, "reset_client", None)
|
|
if reset is not None:
|
|
reset()
|