The confluence engine had nothing to work with. Daily moving averages were the only level source, and they sat 163 to 697 points from price, so every cluster had exactly one member and no alert could ever fire. Two new sources, chosen for having a real following — the engine is a bet that many participants watch the same price, which is what makes a level hold: - Prior day high/low/close, from the last *closed* daily bar so mid-session the levels do not silently switch to today's own developing range. Full daily weight rather than the 0.75 average discount: a traded high is structure, not a derived average. - Session VWAP, anchored to the 18:00 ET open like the daily bars. Institutional execution is benchmarked against it, and zero-volume overnight minutes are skipped rather than dividing by zero. Both are stamped 1d, so they get their own colours to stay distinguishable from the daily averages. Prior-day levels draw as price lines, which span the chart and label the axis instead of relying on bar-index interpolation. VWAP re-prices every minute while a daily average carries hundreds of points and changes once a session, so broadcasting the whole level set on the VWAP cadence would have pushed the entire history every minute. Levels now go out as a delta that clients merge by id. Adding the levels then exposed two defects that had been invisible while nothing could cluster: - Cluster identity was sha1(side + round(center / tolerance)), and tolerance derives from ATR, so it changed every bar. The same zone was continually issued a new id, never matched the cooldown table, and the cooldown did nothing. Identity is now the set of converging levels. - Alert suppression keyed on that identity, so a level drifting in or out of a group read as a new zone. It now suppresses by proximity: two zones within an ATR are the same zone, and the strongest is the one reported. Over six replayed sessions at threshold 28 that is 247 alerts, then 54, then 40; raising the cooldown to 4h — which only affects repeats of the same area, never a genuinely new zone — gives 17 total with a worst session of 9. calibrate_alerts.py now sweeps threshold and cooldown together in one pass, since the threshold turns out to be quantised and nearly useless as a control. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
117 lines
5 KiB
Python
117 lines
5 KiB
Python
import asyncio
|
|
from dataclasses import dataclass, field
|
|
|
|
from app.bars.models import Bar, Timeframe
|
|
from app.bars.aggregator import Aggregator
|
|
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
|
|
|
|
|
|
@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)
|
|
_sent_levels: dict[str, dict] = field(default_factory=dict)
|
|
|
|
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.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)
|
|
self.levels = (
|
|
self.ma_levels
|
|
+ build_prior_day_levels(self.store.get(Timeframe.D1), self.price)
|
|
+ build_vwap_level(minute_bars)
|
|
+ self.manual_lines.levels()
|
|
)
|
|
self.broadcast_level_delta()
|
|
self.rebuild_clusters()
|
|
|
|
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,
|
|
"evaluate_alerts": evaluate_alerts,
|
|
}
|
|
)
|
|
|
|
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")
|