254 lines
11 KiB
Python
254 lines
11 KiB
Python
import asyncio
|
|
import logging
|
|
from dataclasses import dataclass, field, replace
|
|
|
|
from app.analysis.alerts import Alert, AlertEngine
|
|
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
|
|
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})
|
|
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)
|
|
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 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:
|
|
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
|
|
)
|
|
# 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
|
|
return asyncio.create_task(self.stream.run(), name="market-stream")
|