The trendline bug: a line continued past its second anchor at a different slope. Two conventions were fighting, and both were wrong. Lightweight Charts spaces bars evenly however much time separates them — a weekend is forty-nine hours and one bar wide. The renderer extended the line by interpolating between bar indices, which looked straight but disagreed with the server, since price_at() advances per second. Measured on real bars that reached 147 points: the chart drew a level the alerts did not believe in. Making the renderer match price_at() fixed the disagreement and made the visible kick worse, because now the line really did climb an hour's worth of slope across a one-bar maintenance break. Neither convention is what a person means by drawing a line. A trendline advances per bar, so both sides now evaluate in bar space: a new bar_space module the runtime uses to position sloped levels, mirrored by indexAt() in the chart. The line is straight on screen and the alert fires where it is drawn. Also, from testing against the live chart: - A plain click with the trendline tool armed did nothing and left the tool armed, so the next click began a new line — which is how the slope change was first noticed. Click-click and press-drag-release are both supported now, with the rubber band following the cursor between clicks. - Hand-placed levels are armed, fire once, then disarm themselves, and can be re-armed from the sidebar. Verified end to end: created armed, tripped within thirty seconds, disarmed, re-armed. - Layers is collapsible. - The 1h moving averages are gone; only the daily set remains. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
189 lines
8.2 KiB
Python
189 lines
8.2 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.yahoo_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)
|
|
await self.stream.seed(source, Timeframe.H1, self.settings.seed_1h_range)
|
|
await self.stream.seed(source, Timeframe.M1, self.settings.seed_1m_range)
|
|
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")
|