Two things kept the chart quieter than the feed. Higher timeframes only moved once a minute. Tick bars are 1m and the socket filters bar events by the subscriber's timeframe, so on the hourly chart every tick was discarded and only a closed minute passing through the aggregator showed up. They cannot simply be fed to the aggregator — it accumulates with current.v += incoming.v, so the same forming minute re-sent on each tick would add its volume to every higher timeframe again and again. provisional_higher combines the aggregator's committed state with the live minute instead, without mutating it; the next closed minute goes through normally and replaces the result, because the store keys on the bucket timestamp. A test pins the behaviour: five ticks in one minute leave the hour's volume at closed plus live, counted exactly once. Trades known only by their volume were skipped. Level 1 resends only changed fields, so some trades carry a trade stamp and a moved TOTAL_VOLUME with neither LAST_PRICE nor LAST_SIZE. Those now count, with size left at zero rather than guessed from the volume delta — CHART_FUTURES replaces the minute's volume with the exchange's own figure moments later, and two ways of counting the same trades is how double counting starts. Measured: 66 to 74 updates per 90s. The tick throttle drops to 0.25s, which no longer binds. Measured in regular hours the gaps between updates are whole multiples of 1.005s — 2.01, 3.02, 4.03 — which is Schwab conflating LEVEL_ONE_FUTURES to one update per second per symbol. One per second is the source's ceiling, not ours; the longer gaps are seconds in which their feed carried no trade. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
103 lines
3.6 KiB
Python
103 lines
3.6 KiB
Python
import asyncio
|
|
|
|
import pytest
|
|
|
|
from app.analysis.alerts import Alert
|
|
from app.analysis.confluence import Cluster
|
|
from app.analysis.levels import Level, LevelKind, Side
|
|
from app.bars.models import Timeframe
|
|
from app.config import Settings
|
|
from app.runtime import Runtime
|
|
|
|
|
|
def runtime(tmp_path, **overrides) -> Runtime:
|
|
settings = Settings(
|
|
manual_lines_path=tmp_path / "manual_lines.json",
|
|
ntfy_topic=overrides.pop("ntfy_topic", ""),
|
|
**overrides,
|
|
)
|
|
return Runtime(settings)
|
|
|
|
|
|
def alert() -> Alert:
|
|
level = Level(
|
|
"pd:high", LevelKind.HORIZONTAL, Timeframe.D1, Side.RESISTANCE, 16, 1, "PDH",
|
|
100, 5000, 0, None, 0, 100, 100, False, False,
|
|
)
|
|
cluster = Cluster("cl_x", Side.RESISTANCE, 5000, 5000, 5000, 16, [level], 1.0)
|
|
return Alert(cluster, "BEARISH ZONE /ES")
|
|
|
|
|
|
def test_one_engine_serves_every_connection(tmp_path):
|
|
# Previously each WebSocket built its own engine, so reloading the page
|
|
# cleared the cooldown and the same zone alerted again immediately.
|
|
instance = runtime(tmp_path)
|
|
assert instance.alert_engine is not None
|
|
assert instance.alert_engine.min_score == instance.settings.confluence_min_score
|
|
|
|
|
|
def test_alerts_reach_subscribers(tmp_path):
|
|
instance = runtime(tmp_path)
|
|
queue: asyncio.Queue = asyncio.Queue(maxsize=10)
|
|
instance.subscribers.add(queue)
|
|
|
|
async def scenario():
|
|
instance.dispatch_alerts([alert()])
|
|
return queue.get_nowait()
|
|
|
|
event = asyncio.run(scenario())
|
|
assert event["type"] == "alert"
|
|
assert event["message"] == "BEARISH ZONE /ES"
|
|
|
|
|
|
def test_ntfy_failure_does_not_propagate(tmp_path, monkeypatch):
|
|
# send_ntfy used to be awaited inside the WebSocket loop, whose except
|
|
# clause only caught disconnects — so a push outage killed the connection.
|
|
instance = runtime(tmp_path, ntfy_topic="chart-test")
|
|
|
|
async def explode(*args, **kwargs):
|
|
raise RuntimeError("ntfy is down")
|
|
|
|
monkeypatch.setattr("app.runtime.send_ntfy", explode)
|
|
asyncio.run(instance.notify("anything"))
|
|
|
|
|
|
def test_blank_topic_sends_nothing(tmp_path, monkeypatch):
|
|
instance = runtime(tmp_path)
|
|
calls = []
|
|
|
|
async def record(server, topic, message):
|
|
calls.append(topic)
|
|
|
|
monkeypatch.setattr("app.runtime.send_ntfy", record)
|
|
asyncio.run(instance.notify("anything"))
|
|
assert calls == [""] # send_ntfy itself is the one that short-circuits
|
|
|
|
|
|
def test_tick_bars_update_higher_timeframes_without_doubling_volume(tmp_path):
|
|
# The aggregator accumulates with `current.v += incoming.v`, so a forming
|
|
# minute re-sent on every tick would add its volume to each higher
|
|
# timeframe again and again. The provisional path must combine, not
|
|
# accumulate: five ticks in one minute leave the hour's volume equal to the
|
|
# closed minutes plus the live one, exactly once.
|
|
from app.bars.models import Bar
|
|
|
|
instance = runtime(tmp_path)
|
|
|
|
def minute(t, close, volume, closed=True):
|
|
return Bar(tf=Timeframe.M1, t=t, o=close, h=close, l=close, c=close,
|
|
v=volume, closed=closed, symbol="/ES", source="test")
|
|
|
|
base = 1786356000 # top of an hour
|
|
asyncio.run(instance.on_bar(minute(base, 100.0, 10)))
|
|
asyncio.run(instance.on_bar(minute(base + 60, 101.0, 20)))
|
|
settled = [b for b in instance.store.get(Timeframe.H1) if b.t == base][-1].v
|
|
assert settled == 30
|
|
|
|
for _ in range(5):
|
|
asyncio.run(instance.on_bar(minute(base + 120, 102.0, 7, closed=False)))
|
|
|
|
hour = [b for b in instance.store.get(Timeframe.H1) if b.t == base][-1]
|
|
assert hour.v == 37, "the live minute's volume must be added once, not per tick"
|
|
assert hour.c == 102.0
|
|
assert hour.closed is False
|